抱歉,您的浏览器无法访问本站
本页面需要浏览器支持(启用)JavaScript
了解详情 >

DataX 一次同步任务可以压成一条主线:datax.py 组装 JVM 命令,Engine 读取任务 JSON,JobContainer 初始化并切分 Reader/Writer,调度器把 Task 分给若干 TaskGroupContainer,每个 Task 再用 Channel 连接 Reader 和 Writer。

这条链路里最容易混淆的是三个数量:Reader 切出了多少 Task、作业允许多少 Channel、每个 TaskGroup 能同时跑多少 Task。它们相关,但不是一个值。

这篇文章最初有 23 张调试截图,原 OSS 存储桶删除后无法恢复。2026 年整理旧文时,我没有继续保留空图片框,而是按 2020 年 10 月 13 日的 DataX 提交 3c03fedd 重新核对类和方法,并把真正需要看的关系重画成三张调用链图。原始发布时间不变,修订时间单独记录。

DataX 从 Engine 到 Task 执行的主调用链

任务 JSON 先进入 Engine

命令行入口是 bin/datax.py。脚本处理 JVM 参数、日志配置和任务文件路径,最终启动:

1
2
3
4
com.alibaba.datax.core.Engine
-mode standalone
-jobid -1
-job /path/to/job.json

Java 入口在 Engine.entry(args)。它解析 -mode-jobid-job,将任务文件转成 Configuration,再调用 Engine.start(configuration)。Standalone 模式创建的是 JobContainer;TaskGroup 模式才会直接创建 TaskGroupContainer

这里有个很实用的调试办法:不要在 Python 包装脚本里盲猜参数。先把脚本最终生成的 Java 命令完整打印出来,再把 -Ddatax.home、classpath 和三个业务参数放进 IDE Run Configuration。这样断点落在 Engine.entry() 时,配置目录、插件目录和线上执行方式是一致的。

Configuration 贯穿整个执行过程。用户 JSON 只写 job 部分,框架还会合并 core.json、插件元数据和运行参数。后面看到的 Task 配置不是另一套对象,它仍然是这棵配置树切分、克隆后的子树。

JobContainer 管的是作业生命周期

JobContainer.start() 的骨架很稳定:

1
2
3
4
5
6
7
8
preHandle();
init();
prepare();
split();
schedule();
post();
postHandle();
invokeHooks();

我读这段代码时不会逐个方法平均用力。真正决定执行形态的是 init()split()schedule()

init() 根据 Reader/Writer 名称加载 Job 插件,例如 mysqlreadertxtfilewriterLoadUtil 从插件目录读取 plugin.json,建立插件专用类加载器,再实例化 Reader.JobWriter.Job。传给插件的不是完整作业配置,而是对应的 parameter 子树。插件只处理自己的数据源参数,框架配置仍由 Core 持有。

prepare() 允许 Reader 和 Writer 在切分前做一次作业级准备,例如检查表、准备目标目录。它不是每个 Task 都执行一遍。

Reader 决定 Task 粒度,Writer 必须对齐

split() 不是把线程数直接复制成 Task 数。框架先根据速度配置估算 needChannelNumber,再把这个建议数量传给 Reader Job 的 split(adviceNumber)。Reader 可以按分片键、表、文件或数据源连接继续扩展,所以返回数量未必等于 Channel 数。

Writer 随后执行 split(readerTaskConfigs.size())。它必须返回相同数量的配置,框架才能把第 N 个 Reader Task 和第 N 个 Writer Task 合成一个完整 Task。少一个或多一个都不能建立一对一传输链路。

DataX Reader 与 Writer 的 Task 切分和配对

以 MySQL Reader 为例,splitPk 会影响查询范围切分。假设最终生成 25 个 Reader Task,而作业只允许 5 个 Channel,这不代表 25 个 Task 同时跑。它表示调度器有 25 个执行单元,最多并发运行 5 个,完成一批再取下一批。

诊断“为什么只跑了几个并发”时,我会同时记录:

  • job.setting.speed.channelbyterecord
  • JobContainer.adjustChannelNumber() 最终计算出的 Channel 数;
  • Reader split() 实际返回的 Task 数;
  • core.container.taskGroup.channel 规定的单组并发数。

只盯用户 JSON 里的 channel,经常会漏掉总字节预算、单 Channel 预算和 Task 数量对并发的共同约束。限速链路的细节放在另一篇 DataX speed.byte、Channel 与运行时背压分析 里。

schedule() 把 Task 交给 TaskGroup

JobContainer.schedule() 根据总 Channel 数和单 TaskGroup Channel 数计算分组,再调用 Scheduler。Standalone 模式使用进程内调度器;每个 TaskGroup 对应一个 TaskGroupContainer,持有自己的 Task 列表、Channel 配置和重试状态。

TaskGroupContainer.start() 内部维护待运行 Task、正在运行 Task 和失败 Task。循环不会一次性启动所有任务,而是在可用 Channel 出现时创建 TaskExecutor。Task 完成后释放并发名额,失败时根据配置判断是否重试。

这也是排查卡住任务的分界线:

  • Task 还在待运行队列,先查 Channel 名额和前序失败;
  • TaskExecutor 已经启动,查 Reader/Writer 线程和 Channel 等待时间;
  • TaskGroup 没被创建,回到 JobContainer 的切分和调度参数。

把这三个阶段混在一起,日志里看见“任务未运行”就会直接去查数据库,方向很容易错。

一个 Task 里先启动 Writer,再启动 Reader

TaskExecutor 构造时加载对应的 Reader.TaskWriter.Task,同时创建 RecordExchangerdoStart() 会先启动 WriterRunner,再启动 ReaderRunner。这个顺序不是装饰:Writer 先准备消费,Reader 随后生产 Record,可以减少启动瞬间把 Channel 填满的概率。

DataX Reader、Channel 与 Writer 的单 Task 数据流

Reader 插件通过 RecordSender 发送 Record。BufferedRecordExchanger.sendToWriter() 先做字段转换、脏数据处理和批量缓存,再把记录推入 Channel。WriterRunner 把同一个 Channel 暴露成 RecordReceiver,交给 taskWriter.startWrite(recordReceiver) 持续消费。

所以 Reader 和 Writer 并不直接调用对方。它们只依赖 Record 协议和 Channel:

1
2
Reader.Task -> RecordSender -> BufferedRecordExchanger
-> Channel -> RecordReceiver -> Writer.Task

这个边界解释了 DataX 插件为什么容易扩展。新增 Reader 不需要知道目标是 MySQL、HDFS 还是文件;新增 Writer 也不需要知道上游从哪里读。字段类型、脏数据契约和 Channel 容量才是两端共同遵守的东西。

真正有用的运行证据

源码读完以后,线上诊断仍然要落回运行证据。我通常保留下面几组值:

1
2
3
4
5
6
任务配置摘要:reader / writer / speed / errorLimit
切分结果:readerTaskCount / writerTaskCount
调度结果:channelCount / taskGroupCount / tasksPerGroup
传输统计:byteSpeed / recordSpeed / totalReadRecords
等待统计:waitReaderTime / waitWriterTime
失败位置:Job prepare / split / Task Reader / Channel / Task Writer

waitWriterTime 高不等于 Writer 在等。它记录 Reader 向 Channel 推送时等待空间的时间,往往说明 Writer 或目标端较慢;waitReaderTime 高则通常是 Writer 等待数据。指标名要顺着采集位置解释,不能按字面猜。

还有一条我会坚持:先固定数据集、Reader SQL、Writer 批次和运行时段,再只改一个并发或限速参数。没有对照组的“调大后快了”很难区分缓存、数据分布和目标端负载。

旧图恢复后的边界

原文的大部分截图只是 IDE 中的一段方法或展开后的 JSON。它们在 2020 年能证明我确实沿着代码调试过,但对今天的读者帮助有限,而且原图已经不可验证。这次保留了启动参数、配置形态和核心方法,删除无法恢复的截图式证据,改成可版本化的 PlantUML。

没有恢复的内容也明确留下边界:本文没有给出生产吞吐数字,不把某组本地调试结果包装成通用性能结论。需要复现时,应 checkout 同一提交,使用自己的 Reader/Writer 配置跑断点和日志。

源码对照